Java 响应式编程 - Reactive Streams 规范

在现代高并发、高吞吐量的互联网应用中,传统的一个请求对应一个线程模型很容易在 I/O 阻塞(如数据库查询、远程调用)时导致线程资源耗尽。为了解决这个问题,响应式编程(Reactive Programming)应运而生。而 Reactive Streams(响应式流)规范,则是整个 Java 生态响应式编程的基石。


什么是 Reactive Streams 规范

Reactive Streams 是一套面向流式数据处理、具有 非阻塞背压(Non-blocking Backpressure)的异步组件标准。简单来说,它定义了一套 API 规范,让不同的框架(如 Spring WebFlux、RxJava、Akka Streams)可以互相兼容,共同解决 “上游生产数据太快,下游消费太慢,导致内存溢出(OOM)” 的核心痛点。

Reactive Streams 的核心特征就是两个字——背压。在传统的流处理中,如果上游(Publisher)发送数据的速度远快于下游(Subscriber)处理的速度,下游的内存缓冲区就会被撑爆。背压机制允许下游订阅者显式地告诉上游:“我现在只能处理 3 个元素,请不要发多了。” 从而实现了流量控制,保证了系统的稳定性。

Reactive Streams 官网 reactive-streams.org。它的 Java 实现非常精简,在 Java 9 之中被直接引入到了 java.util.concurrent.Flow 类中。

该规范极其精简,核心只有 4 个接口(四大核心组件):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
// 1. 发布者:数据的源头,负责生产数据
public interface Publisher<T> {
public void subscribe(Subscriber<? super T> s);
}

// 2. 订阅者:数据的终点,负责消费数据
public interface Subscriber<T> {
public void onSubscribe(Subscription s); // 订阅成功时的回调
public void onNext(T t); // 接收到新数据时的回调
public void onError(Throwable t); // 发生错误时的回调
public void onComplete(); // 数据流正常结束时的回调
}

// 3. 订阅关系:连接发布者和订阅者的纽带,背压的核心(最核心的灵魂)
public interface Subscription {
public void request(long n); // 下游向上游请求 n 个数据(背压控制)
public void cancel(); // 下游取消订阅
}

// 4. 处理器:既是发布者又是订阅者,用于对数据进行转换(中间操作)
public interface Processor<T, R> extends Subscriber<T>, Publisher<R> {
}


Reactive Streams 的原理和流程

并发模型

Reactive Streams 规范本身没有绑定任何特定的线程模型,它是一套协议。我们这里说的其原理和工作流程对应到它的落地到具体实现,比如 Project Reactor 和 Spring WebFlux 上。

第一,在多线程并发下,下游可能在主线程请求数据,上游可能在定时器线程产生数据。为了保证高并发下的绝对安全且不阻塞,Subscription 内部通常使用 CAS 乐观锁和原子变量(如 AtomicLong、AtomicReference)来维护两个核心状态:

  • requested(请求计数器):下游请求的剩余安全配额。每次下游 request(n),利用 CAS 执行 requested += n;每次上游 onNext(),执行 requested–。当 requested == 0 时,上游必须彻底停止发送,原地待命。
  • cancelled(取消标志):布尔原子量。一旦变为 true,后续所有 onNext 直接被丢弃并清理内存。

第二,在传统多线程编程中,我们需要自己处理 ExecutorService、CountDownLatch、线程间共享队列的可见性、死锁等问题。Reactive Streams 框架(以 Project Reactor 为例)通过将 “线程切换” 抽象为声明式的操作符,彻底屏蔽了这些复杂性:

  • publishOn(Scheduler):改变下游执行的线程池。
  • subscribeOn(Scheduler):改变上游数据源加载时执行的线程池。

第三,当你写下 flux.publishOn(Schedulers.boundedElastic()).map(…) 时,Reactor 底层实际上利用了异步环形队列(Ring Buffer / Queue)和事件驱动回调:

  • 无锁队列隔离:publishOn 内部会偷偷塞入一个 Processor 或一个高性能的单生产者-单消费者(SPSC)异步无锁队列。
  • 线程切换点:上游线程(假设是 A 线程)生产了数据,它不直接调用下游的 onNext,而是把数据塞进这个队列,并向目标线程池(BoundedElastic 线程池中的 B 线程)提交一个执行任务。
  • 消费者线程接管:B 线程苏醒后,从队列中源源不断地取出数据,然后高效率地执行下游的 map 和 subscribe 逻辑。

这一切都被封装在操作符内部,开发者看到的只是优雅的链式调用。在应用运行时,底层的线程开多少、怎么开,完全取决于宿主服务器(如 Netty)和调度器(Schedulers)的配置。我们以 Spring WebFlux + Netty + Reactor 这一黄金组合来拆解:

① 核心网络 I/O 线程池:Netty EventLoopGroup:当你启动一个 WebFlux 应用时,底层默认运行着一个 Netty 服务器。

  • 开了多少线程?:默认是 CPU 核心数 x 2​(通常被称为 reactor-http-nio-X 线程)。
  • 怎么工作?:这些是极度宝贵的非阻塞 Worker 线程。它们利用操作系统的多路复用(如 Linux 的 epoll),一个线程就能同时管理成千上万个网络连接的 Socket 读写。
  • 致命禁忌:绝对不能在这些线程里执行任何阻塞操作(如 JDBC、Thread.sleep())。一旦阻塞,Netty 的事件循环就会卡死,整个服务器的吞吐量会瞬间雪崩。

② Reactor 的内置调度器(线程池类型)

为了处理不同类型的业务,Reactor 提供了几类开箱即用的线程池(Schedulers):


核心工作流程

1
2
3
4
5
6
7
8
Publisher                Subscriber
| |
|---- subscribe(s) ------>|
|<--- onSubscribe(sub) ---| (传递 Subscription)
| |
|<--- request(n) ---------| (下游主动索要 n 个数据)
|---- onNext(item) ------>| (上游发送数据,最多 n 个)
|---- onComplete()/Error->| (结束或异常)
  • 建立连接 (subscribe):Subscriber 调用 Publisher.subscribe(Subscriber)。
  • 初始化并传递令牌 (onSubscribe):Publisher 创建一个 Subscription 对象(常命名为 ScalarSubscription 或 FluxHide.SuppressSubscription),通过 Subscriber.onSubscribe(Subscription) 传给下游。此时,上游不会发送任何数据。
  • 下游按需索要 (request):下游 Subscriber 根据自身处理能力,调用 Subscription.request(n)。这一步通常是异步的。
  • 上游生产并推送 (onNext):上游收到 request(n) 信号后,开始生产数据,并调用 Subscriber.onNext(T) 传递数据。上游发送的数据量绝对不能超过 n​。
  • 循环或结束 (onComplete / onError):下游消费完这 $n$ 个数据后,如果还需要,就继续调用 request(m);如果数据发完了,上游调用 onComplete() 关闭流。


一个完整的请求

我们再来看一个 WebFlux 收到请求并调用 R2DBC(Reactive Relational Database Connectivity)的底层真实全链路线程流转轨迹:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
[用户浏览器] 
│ (发起 HTTP 请求)

[Netty 线程池: reactor-http-nio-1]
│ 1. 负责接收网卡数据,解析 HTTP 协议
│ 2. 调用 WebFlux 路由,触发 Flux/Mono 链条
│ 3. 向响应式数据库驱动发送 SQL 请求(非阻塞 Netty Channel)
│ 4. 立刻释放线程!(reactor-http-nio-1 去接待其他用户了,此时该线程完全空闲)

[操作系统内核 / 数据库异步等待] ─── (此时没有任何 Java 线程在阻塞等待)

▼ (数据库查询完毕,网卡收到数据)
[R2DBC 驱动线程 / Netty 读线程: reactor-http-nio-2]
│ 1. 驱动收到内核通知,数据已就绪
│ 2. 触发 Reactive Streams 的对游信号:Subscription.request(n)
│ 3. 触发下游的 Subscriber.onNext(data)

[WebFlux 核心线程: reactor-http-nio-2]
│ 1. 执行你的业务代码 map() / filter()
│ 2. 将结果序列化为 JSON,通过网卡非阻塞写回用户

你可能会有疑惑:整个过程中,查询数据库这一步不是阻塞的吗?如果是异步调用,数据库返回数据又是怎么返回 前端的呢?绝大多数从传统 Servlet(Spring MVC)转到响应式编程的开发者,都会在这里卡壳:数据库查询怎么可能不阻塞?哪怕用非阻塞,数据回来后又是怎么神不知鬼不觉地塞回前端的?这里隐藏着两个底层的关键技术:R2DBC(驱动层的革新) 和 Netty Channel Pipeline(网络层的复用)。

传统的 JDBC(如 mysql-connector-java)之所以是阻塞的,是因为它的底层通信使用了传统的 BIO(Blocking I/O)。当 Java 线程执行 resultSet.next() 时,线程会发出一个 TCP 请求,然后原地死等数据库服务器的响应。此时操作系统会将该线程挂起(Block),直至收到数据。

而响应式编程使用的 R2DBC,其底层彻底抛弃了 BIO,改用 NIO(Non-blocking I/O,通常基于 Netty)。底层真实的执行过程是:

  • 发送请求后立即撒手:当你的业务代码调用数据库查询时,Java 线程(Netty Worker 线程)将 SQL 语句打包成一个 TCP 数据包,通过操作系统的非阻塞 Socket 发送出去。
  • 不等待,直接返回:发送完毕后,该线程绝对不会在原地等待数据库返回结果。它在驱动层注册一个 “数据准备就绪” 的回调事件(Event Listener),然后这个线程就立刻结束当前任务,去高高兴兴地接待下一个用户的 HTTP 请求了。
  • 内核接管:此时,没有任何 Java 线程在为这个数据库查询买单。只有操作系统的内核(如 Linux 的 epoll)在盯着网卡。

既然发送请求的线程已经跑去干别的事了,那么当数据库执行完 SQL,把数据顺着网络线缆丢回来的时候,数据是怎么一步步回到前端浏览器的呢?这里全靠操作系统的事件通知机制和 Reactive Streams 的链式订阅关系来完成接力。整个过程就像一个完美的多米诺骨牌效应:

  • 骨牌 1:网卡收到数据,内核唤醒 Netty。数据库把数据计算好后,通过 TCP 回传。服务器的网卡收到数据包,引发一个硬件中断。Linux 内核发现这个 Socket 有数据进来了,于是通知正在执行 epoll_wait 的 Netty 线程池(可能是刚才那个线程,也可能是同池子里的另一个空闲 Worker 线程,我们假设是 reactor-http-nio-2)。
  • 骨牌 2:驱动程序触发 Reactive Streams 信号。reactor-http-nio-2 线程领到任务,开始读取网卡缓冲区里的数据库原始字节流。R2DBC 驱动将这些字节解析成 Java 对象。此时,驱动程序作为 Publisher(发布者),开始向下游触发 Reactive Streams 标准信号:
    • 它调用下游的 Subscription.request(n) 看需不需要数据。
    • 确认需要后,调用 Subscriber.onNext(data) 将转换好的数据行(Row)推送出去。
  • 骨牌 3:业务操作符链条被 “串路激活”。数据一旦进入 onNext,你之前写的那些声明式代码(如 .map(user -> …), .filter(…))就开始在 reactor-http-nio-2 线程上流水线式地执行。这些代码并不是在请求时执行的,而是在数据回来的这一刻,被数据流“冲刷”激活的。
  • 骨牌 4:Netty 写入网络缓冲区,直达前端。当数据流一路走到最下游(WebFlux 框架层)时,WebFlux 的底层订阅者会把这些数据对象序列化为 JSON 字符串,然后直接调用 Netty 的 channel.writeAndFlush()。 Netty 再次利用非阻塞 I/O,把 JSON 字节写入到与前端浏览器建立的那个 Socket 管道(Channel)中。网卡将数据发送出去,前端浏览器收到响应,渲染页面。

通过这种 “事件驱动、收到通知才干活” 的方式,哪怕数据库查询需要耗时 3 秒,Java 服务器在这 3 秒内也不需要耗费任何一个线程去死等,从而释放了恐怖的并发吞吐能力。


两个具体的测试案例

基础案例

以下是一个只包含发布者和订阅者的最简单的案例:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
import java.io.IOException;
import java.util.concurrent.Flow;
import java.util.concurrent.SubmissionPublisher;

public class ReactiveDemo {

public static void main(String[] args) throws IOException {
// 1. 创建发布者 (JDK 自带的 SubmissionPublisher 实现了 Publisher)
SubmissionPublisher<String> publisher = new SubmissionPublisher<>();

// 2. 创建订阅者
Flow.Subscriber<String> subscriber = new Flow.Subscriber<>() {
private Flow.Subscription subscription;

@Override
public void onSubscribe(Flow.Subscription subscription) {
this.subscription = subscription;
// 订阅成功后,先向游请求 1 个数据
this.subscription.request(1);
}

@Override
public void onNext(String item) {
System.out.println("【接收数据】: " + item);
// 消费完后,继续请求 1 个数据(背压的体现)
this.subscription.request(1);
}

@Override
public void onError(Throwable e) {
e.printStackTrace();
}

@Override
public void onComplete() {
System.out.println("【数据流处理完成】");
}
};

// 3. 建立订阅关系
publisher.subscribe(subscriber);

// 4. 发布数据
System.out.println("【开始发布数据】");
publisher.submit("Hello");
publisher.submit("Reactive");
publisher.submit("Streams");

// 关闭发布者
publisher.close();

// 主线程等待异步执行结束
System.in.read();
}
}


包含 Processor 的案例

Processor(处理者)同时继承了 Subscriber 和 Publisher,这意味着它扮演了一个“双面间谍” 的角色:

  • 对于上游,它是一个订阅者,负责接收原始数据;
  • 对于下游,它是一个发布者,负责把加工、过滤、转换后的数据推送出去。

下面的案例中:Processor 拦截上游的原始订单,自动过滤掉 100 元以下的“小额订单”,并将剩余的高端订单包装成 “VIP专属订单” 再投递给下游的风控系统。这个案例完整展示了底层的链式接力机制:

Publisher (数据源) —–> Processor (加工过滤) —–> Subscriber (最终消费)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;

import java.io.IOException;
import java.util.Arrays;
import java.util.List;
import java.util.concurrent.Flow;
import java.util.concurrent.atomic.AtomicInteger;

public class ReactiveDemo {

@Data
@AllArgsConstructor
@NoArgsConstructor
static class Order {
String orderId;
int amount;
}

static class OrderFilterProcessor implements Flow.Processor<Order, Order> {

// 作为发布者:它需要持有它的下游订阅者
private Flow.Subscriber<? super Order> actualSubscriber;
// 作为订阅者:它需要持有它的上游订阅令牌
private Flow.Subscription upstreamSubscription;

@Override
public void subscribe(Flow.Subscriber<? super Order> subscriber) {
// 记录下真正的下游是谁
this.actualSubscriber = subscriber;
// 规范要求:当下游订阅此 Processor 时,必须把控制令牌传给下游。这里我们通过一个桥接的匿名 Subscription 把控制权转交出去
this.actualSubscriber.onSubscribe(new Flow.Subscription() {
@Override
public void request(long n) {
System.out.printf("[Processor] 收到下游拉取信号 request(%d),直接透传转发给上游...\n", n);
// 核心背压透传:下游要多少,Processor 就去跟上游要多少
upstreamSubscription.request(n);
}

@Override
public void cancel() {
System.out.println("[Processor] 收到下游取消信号,通知上游取消...");
upstreamSubscription.cancel();
}
});
}

@Override
public void onSubscribe(Flow.Subscription subscription) {
// 记录上游给它的令牌
this.upstreamSubscription = subscription;
}

@Override
public void onNext(Order item) {
System.out.println("[Processor 拦截加工] 收到订单: " + item.orderId + " (" + item.amount + "元)");

// 核心加工过滤逻辑:金额小于 100 元的不发给下游
if (item.amount >= 100) {
System.out.println("[Processor 过滤通过] 金额符合 VIP 标,推送给下游风控系统...");
// 把加工或筛选过后的数据送给真正的下游
this.actualSubscriber.onNext(item);
} else {
System.out.println("[Processor 拦截剔除] 金额太小 (<100元),直接抛弃。");
// 关键点:由于抛弃了一个元素,下游没有拿到数据,为了不破坏背压数量,必须主动再去跟上游要 1 个补上!
this.upstreamSubscription.request(1);
}
}

@Override
public void onError(Throwable throwable) {
this.actualSubscriber.onError(throwable); // 异常向下透传
}

@Override
public void onComplete() {
System.out.println("[Processor] 上游数据发完了,通知下游结束。");
this.actualSubscriber.onComplete(); // 结束信号向下透传
}
}

// 最终订阅者
static class RiskSubscriber implements Flow.Subscriber<Order> {
private Flow.Subscription sub;

@Override
public void onSubscribe(Flow.Subscription s) {
this.sub = s;
System.out.println("[Subscriber] 成功订阅,请求 1 个高端订单...");
this.sub.request(1); // 每次只要 1 个
}

@Override
public void onNext(Order item) {
System.out.println("[Subscriber 最终消费] 收到优质大单并完成风控: " + item);
System.out.println("[Subscriber] 处理完当前订单,继续请求下 1 个...");
this.sub.request(1);
}

@Override
public void onError(Throwable t) {}

@Override
public void onComplete() {
System.out.println("[Subscriber] 完美的流式结束!");
}
}

// 自定义的简单 Subscription
static class SimpleSubscription implements Flow.Subscription {
private final Flow.Subscriber<? super Order> sub;
private final List<Order> data;
private final AtomicInteger index = new AtomicInteger(0);

public SimpleSubscription(Flow.Subscriber<? super Order> sub, List<Order> data) {
this.sub = sub;
this.data = data;
}

@Override
public void request(long n) {
for (int i = 0; i < n; i++) {
int curr = index.getAndIncrement();
if (curr < data.size()) {
sub.onNext(data.get(curr));
} else {
sub.onComplete();
break;
}
}
}

@Override
public void cancel() {}
}

// 原始发布者
static class OrderPublisher implements Flow.Publisher<Order> {
private final List<Order> db = Arrays.asList(
new Order("A001", 150), // 通过
new Order("A002", 50), // 过滤删除
new Order("A003", 300) // 通过
);

@Override
public void subscribe(Flow.Subscriber<? super Order> s) {
s.onSubscribe(new SimpleSubscription(s, db));
}
}


public static void main(String[] args) throws IOException {
OrderPublisher publisher = new OrderPublisher();
OrderFilterProcessor processor = new OrderFilterProcessor();
RiskSubscriber subscriber = new RiskSubscriber();

// 核心链条串联:
// 1. 处理器订阅原始发布者
publisher.subscribe(processor);
// 2. 最终订阅者订阅处理器
processor.subscribe(subscriber);

System.in.read();
}
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
[Subscriber] 成功订阅,请求 1 个高端订单...
[Processor] 收到下游拉取信号 request(1),直接透传转发给上游...
[Processor 拦截加工] 收到订单: A001 (150元)
[Processor 过滤通过] 金额符合 VIP 标,推送给下游风控系统...
[Subscriber 最终消费] 收到优质大单并完成风控: ReactiveDemo.Order(orderId=A001, amount=150)
[Subscriber] 处理完当前订单,继续请求下 1 个...
[Processor] 收到下游拉取信号 request(1),直接透传转发给上游...
[Processor 拦截加工] 收到订单: A002 (50元)
[Processor 拦截剔除] 金额太小 (<100元),直接抛弃。
[Processor 拦截加工] 收到订单: A003 (300元)
[Processor 过滤通过] 金额符合 VIP 标,推送给下游风控系统...
[Subscriber 最终消费] 收到优质大单并完成风控: ReactiveDemo.Order(orderId=A003, amount=300)
[Subscriber] 处理完当前订单,继续请求下 1 个...
[Processor] 收到下游拉取信号 request(1),直接透传转发给上游...
[Processor] 上游数据发完了,通知下游结束。
[Subscriber] 完美的流式结束!


关于响应式编程的学习

响应式编程(Reactive Programming)也许你刚开始觉得很晦涩,是因为它彻底颠覆了我们写了十几年的命令式编程(Imperative Programming)思维。过去我们习惯了先 A 再 B 卡住等结果,而响应式关注的是数据怎么流动,怎么对变化做出响应。为了让你彻底看清,我们跳出复杂的 API,用一条 “由浅入深、层层递进” 的黄金主线来帮你把零散的知识点串联起来。

  • 第一站:思维重构(从拉模型到推模型)
    • 传统思维(拉模型 - Pull):你作为老板(主线程),让员工 A 去数据库拿数据。A 去了,你就在办公室坐着死等(阻塞),A 回来之前你什么都干不了。
    • 响应式思维(推模型 - Push):你(主线程)不再等 A,而是构建一条生产流水线/自来水管。你跟 A 说:“你拿到数据后,顺着这条管道丢进来。” 交代完之后,你立刻去干别的事。当数据真正顺着管道流下来的时候,流水线上的各个节点自动触发加工。
    • 响应式编程 = 声明式数据流(Data Stream) + 变化传递(Propagation of change)。你在写代码时,不是在下达执行命令,而是在“搭管道”。
  • 第二站:统一标准(Reactive Streams 规范)
    • 当你理解了管道思维,就会发现一个问题:上游放水太快,下游池子太小,水会溢出来(也就是内存溢出 OOM)。于是,Java 制定了 Reactive Streams 规范。
    • 规范极其简单,它用 4 个接口规定了管道两端和中间的玩法:
      • Publisher(发布者 / 水龙头):负责生产数据。
      • Subscriber(订阅者 / 盛水的水桶):负责消费数据。
      • Subscription(订阅关系 / 水阀开关)最核心的灵魂! 下游调用 request(n) 告诉上游:“我现在有空,给我放 $n$ 滴水。” 上游放完 $n$ 滴立刻停下,等下次通知。这就是背压机制。
      • Processor(处理器 / 过滤器):既是水桶又是水龙头,放在中间净水。
      • 规范不包含任何复杂的线程池或高深技术,它只是通过 request(n) 这一小步,完美解决了 “上游太快饿死/撑死下游” 的千古难题。
  • 第三站:乐高组装(Project Reactor 核心)
    • 规范只是设计图纸,无法直接在生产中使用。Pivotal 团队(Spring 亲爹)根据这个规范写出了工业级实现——Project Reactor。它把晦涩的接口变成了如同乐高积木般好玩的工具。它提供了两个核心管道工:
      • Mono\<T>:只会流出 0 或 1 个元素的特制管道。
      • Flux\<T>:可以流出 0 到 N 个无穷元素的标准管道。
    • 它最大的魅力在于提供了几百个操作符(Operators),就像各式各样的管道配件:
      • filter():过滤掉不合格的水(数据)。
      • map():把清水变成纯净水(数据类型转换)
      • publishOn() / subscribeOn():魔法切换器。水流到这里,自动切换到另一个动力泵(线程池)继续推,开发者完全不用操心底层的 Thread 和 Lock。
      • 完全不要去死记几百个操作符!记住它们只是管道配件,随用随查。你只需要用链式调用把 Flux.just(1, 2).map(…).filter(…) 串起来即可。
  • 第四站:实战落地(Spring WebFlux)
    • 当你学会了组装管道,最后一步就是用它来重构你的 Web 应用。这就是 Spring WebFlux。它彻底抛弃了 Tomcat 那种“一个请求死死绑定一个线程”的模式,改用 Netty(基于事件驱动的非阻塞服务器)。
      • 传统 Tomcat:1000 个并发请求,必须要 1000 个物理线程。如果数据库很慢,1000 个线程全部卡死。
      • WebFlux + Netty:可能一共就 4 个核心线程。遇到 I/O 阻塞(比如查数据库),这 4 个线程“把请求往操作系统的异步队列里一丢”,立刻转头去接第 5 个请求。等到数据库数据回来了,操作系统发个通知,4 个线程中的某一个再回来把数据顺着管道推给前端。
      • WebFlux 的并发能力不是靠增加线程实现的,而是靠 “不让线程闲着” 实现的。
  • 避坑建议:
    • 忘记 ThreadLocal:不要再奢望用 ThreadLocal 传递登录信息或日志 MDC 令牌了,因为响应式下,你的代码可能前 1 秒在线程 A 执行,后 1 秒就在线程 B 执行了。要习惯使用 Reactor 的 Context。
    • 绝对不要在管道里“下毒”:在 Flux/Mono 的链条里,绝对绝对不能调用阻塞方法(比如老版 HttpClient.execute()、Thread.sleep() 或传统的 JDBC)。只要有一个阻塞,整个响应式线程池就会瘫痪。
    • 记住“不订阅,什么都不会发生”:比如 Mono\<User> user = r2dbcRepository.findById(1L) ,这行代码执行了,但由于最后没有 subscribe(),数据库连查都不会查,它只是一张设计图!

按照 “重构思维 $\rightarrow$ 搞懂背压接口 $\rightarrow$ 熟悉 Flux/Mono 链式操作 $\rightarrow$ 结合 WebFlux” 这条主线去学,你会发现那些晦涩的 API 背后,其实是一套极其优雅、有条不紊的现代化工业流水线。